Skip to content

[fix](streaming-job) Optimize snapshot offset persistence - #66238

Open
JNSimba wants to merge 2 commits into
apache:masterfrom
JNSimba:fix/streaming-snapshot-offset-persist
Open

[fix](streaming-job) Optimize snapshot offset persistence#66238
JNSimba wants to merge 2 commits into
apache:masterfrom
JNSimba:fix/streaming-snapshot-offset-persist

Conversation

@JNSimba

@JNSimba JNSimba commented Jul 29, 2026

Copy link
Copy Markdown
Member

What problem does this PR solve?

Problem Summary:

Snapshot tasks persisted the complete offset provider after every split commit, producing frequent edit logs while the snapshot state kept growing. The previous default snapshot split size also produced too many small snapshot tasks. Natural terminal transitions needed explicit persistence so finished jobs replay with their terminal metadata.

This change:

  • Limits snapshot offset edit logs with a mutable 300-second interval while keeping binlog offset commits immediate.
  • Increases the default snapshot split size from 8,096 to 40,960 rows, while preserving explicitly configured values.
  • Removes obsolete snapshot recovery state after the first committed binlog offset for FROM-TO and CDC stream TVF jobs.
  • Persists natural FINISHED transitions and restores terminal replay fields and callbacks consistently.
  • Adds focused unit coverage and a MySQL snapshot-only FE restart regression case.

Release note

None

Check List (For Author)

  • Test

    • Regression test (added; execution pending)
    • Unit Test (added; execution pending)
    • Manual test
    • No need to test or manual test.
  • Behavior changed:

    • No.
    • Yes. Snapshot offset edit-log persistence is rate-limited, the default snapshot split size is increased to 40,960 rows, and binlog offsets remain immediately persisted.
  • Does this need documentation?

    • No.
    • Yes.

Check List (For Reviewer who merge this PR)

  • Confirm the release note
  • Confirm test cases
  • Confirm document
  • Add branch pick label

@hello-stephen

Copy link
Copy Markdown
Contributor

Thank you for your contribution to Apache Doris.
Don't know what should be done next? See How to process your PR.

Please clearly describe your PR:

  1. What problem was fixed (it's best to include specific error reporting information). How it was fixed.
  2. Which behaviors were modified. What was the previous behavior, what is it now, why was it modified, and what possible impacts might there be.
  3. What features were added. Why was this function added?
  4. Which code was refactored and why was this part of the code refactored?
  5. Which functions were optimized and what is the difference before and after the optimization?

@JNSimba

JNSimba commented Jul 29, 2026

Copy link
Copy Markdown
Member Author

/review

@JNSimba

JNSimba commented Jul 29, 2026

Copy link
Copy Markdown
Member Author

run buildall

@hello-stephen

Copy link
Copy Markdown
Contributor
TPC-H: Total hot run time: 29504 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpch-tools
Tpch sf100 test result on commit fa4970577ded6215f08aea5bad215feb85a2840f, data reload: false

------ Round 1 ----------------------------------
============================================
q1	17625	4129	4015	4015
q2	2375	326	196	196
q3	10279	1456	838	838
q4	4700	475	342	342
q5	7861	858	586	586
q6	203	170	138	138
q7	789	821	631	631
q8	10102	1596	1584	1584
q9	7043	4298	4317	4298
q10	7356	1734	1447	1447
q11	525	351	334	334
q12	743	582	452	452
q13	18254	3390	2755	2755
q14	270	267	238	238
q15	q16	789	786	708	708
q17	981	997	950	950
q18	7136	5822	5647	5647
q19	1528	1309	1033	1033
q20	816	672	600	600
q21	5756	2737	2415	2415
q22	436	358	297	297
Total cold run time: 105567 ms
Total hot run time: 29504 ms

----- Round 2, with runtime_filter_mode=off -----
============================================
q1	4449	4321	4312	4312
q2	295	321	214	214
q3	4566	4964	4436	4436
q4	2055	2143	1381	1381
q5	4399	4257	4741	4257
q6	256	196	145	145
q7	2127	1846	1619	1619
q8	2450	2084	2099	2084
q9	7783	7773	7763	7763
q10	4730	4699	4448	4448
q11	584	402	389	389
q12	740	761	543	543
q13	3336	3626	2982	2982
q14	304	316	288	288
q15	q16	713	717	649	649
q17	1339	1326	1316	1316
q18	7982	7366	6869	6869
q19	1137	1081	1090	1081
q20	2234	2229	1935	1935
q21	5205	4534	4451	4451
q22	509	486	435	435
Total cold run time: 57193 ms
Total hot run time: 51597 ms

@hello-stephen

Copy link
Copy Markdown
Contributor
TPC-DS: Total hot run time: 178444 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpcds-tools
TPC-DS sf100 test result on commit fa4970577ded6215f08aea5bad215feb85a2840f, data reload: false

query5	4361	622	495	495
query6	470	240	202	202
query7	5043	608	333	333
query8	350	187	170	170
query9	8776	4103	4101	4101
query10	455	362	303	303
query11	5929	2345	2133	2133
query12	162	107	108	107
query13	1267	583	426	426
query14	6255	5298	4938	4938
query14_1	4276	4249	4248	4248
query15	213	204	178	178
query16	1070	486	395	395
query17	1137	704	592	592
query18	2460	474	354	354
query19	219	195	159	159
query20	116	112	106	106
query21	237	161	141	141
query22	13650	13588	13359	13359
query23	17435	16490	16168	16168
query23_1	16322	16311	16392	16311
query24	7610	1748	1290	1290
query24_1	1311	1323	1302	1302
query25	579	474	391	391
query26	1360	368	220	220
query27	2543	641	405	405
query28	4392	2006	2015	2006
query29	1076	645	489	489
query30	345	272	236	236
query31	1117	1104	1007	1007
query32	111	64	62	62
query33	549	347	269	269
query34	1225	1181	648	648
query35	778	811	682	682
query36	1020	1036	903	903
query37	158	108	109	108
query38	1876	1718	1650	1650
query39	873	897	906	897
query39_1	833	869	848	848
query40	251	171	143	143
query41	78	60	62	60
query42	94	93	93	93
query43	322	321	280	280
query44	1435	779	760	760
query45	200	183	174	174
query46	1131	1205	751	751
query47	2144	2144	2022	2022
query48	396	450	301	301
query49	608	440	324	324
query50	1105	442	359	359
query51	10457	10736	10766	10736
query52	93	91	80	80
query53	256	281	197	197
query54	290	241	217	217
query55	77	71	65	65
query56	313	292	289	289
query57	1320	1287	1172	1172
query58	288	267	251	251
query59	1579	1653	1486	1486
query60	313	274	258	258
query61	150	155	147	147
query62	543	502	438	438
query63	242	207	194	194
query64	2696	1034	824	824
query65	4755	4657	4684	4657
query66	1797	512	405	405
query67	29400	29343	29249	29249
query68	3203	1612	1042	1042
query69	419	305	264	264
query70	895	833	782	782
query71	375	357	325	325
query72	3083	2631	2412	2412
query73	861	853	442	442
query74	5087	4932	4714	4714
query75	2520	2502	2145	2145
query76	2333	1208	789	789
query77	359	388	289	289
query78	11886	12010	11431	11431
query79	2233	1201	758	758
query80	1528	583	467	467
query81	560	330	285	285
query82	593	158	121	121
query83	368	326	313	313
query84	275	155	128	128
query85	986	638	537	537
query86	429	238	229	229
query87	1821	1830	1770	1770
query88	3759	2815	2795	2795
query89	434	381	336	336
query90	1888	201	201	201
query91	206	192	162	162
query92	64	64	57	57
query93	1808	1530	999	999
query94	714	368	316	316
query95	797	512	491	491
query96	1133	771	350	350
query97	2643	2614	2480	2480
query98	213	208	205	205
query99	1121	1120	1031	1031
Total cold run time: 265474 ms
Total hot run time: 178444 ms

@hello-stephen

Copy link
Copy Markdown
Contributor
ClickBench: Total hot run time: 25.17 s
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/clickbench-tools
ClickBench test result on commit fa4970577ded6215f08aea5bad215feb85a2840f, data reload: false

query1	0.00	0.00	0.00
query2	0.10	0.04	0.04
query3	0.26	0.14	0.14
query4	1.62	0.15	0.14
query5	0.24	0.22	0.22
query6	1.24	1.10	1.11
query7	0.04	0.01	0.01
query8	0.05	0.04	0.04
query9	0.38	0.31	0.32
query10	0.57	0.59	0.56
query11	0.20	0.13	0.13
query12	0.18	0.15	0.15
query13	0.49	0.49	0.50
query14	1.02	1.03	1.02
query15	0.62	0.61	0.61
query16	0.30	0.32	0.33
query17	1.12	1.09	1.13
query18	0.22	0.21	0.22
query19	2.06	1.92	1.95
query20	0.01	0.01	0.01
query21	15.44	0.19	0.14
query22	4.94	0.05	0.06
query23	16.13	0.29	0.12
query24	3.15	0.41	0.33
query25	0.11	0.06	0.05
query26	0.73	0.21	0.16
query27	0.05	0.03	0.04
query28	3.50	1.00	0.51
query29	12.48	4.20	3.30
query30	0.27	0.15	0.15
query31	2.76	0.60	0.32
query32	3.23	0.60	0.49
query33	3.14	3.21	3.15
query34	15.57	4.29	3.55
query35	3.52	3.51	3.52
query36	0.54	0.43	0.45
query37	0.09	0.07	0.06
query38	0.05	0.04	0.04
query39	0.04	0.03	0.03
query40	0.18	0.16	0.15
query41	0.14	0.03	0.03
query42	0.05	0.03	0.03
query43	0.04	0.03	0.04
Total cold run time: 96.87 s
Total hot run time: 25.17 s

@hello-stephen

Copy link
Copy Markdown
Contributor

FE UT Coverage Report

Increment line coverage 81.58% (31/38) 🎉
Increment coverage report
Complete coverage report

@github-actions github-actions Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed the full authoritative change set. I found one merge-blocking regression-test authoring issue; I did not find a PR-introduced production correctness defect.

Part 1.3 checkpoint conclusions:

  1. Goal and proof: The change reduces snapshot-offset EditLog volume, compacts obsolete snapshot recovery state after the binlog transition, durably records natural FINISHED state, and supplies a 40960 default snapshot split size. The production implementation accomplishes those goals. Focused FE tests exercise persistence timing, binlog immediacy, state compaction, replayed terminal fields/callback cleanup, and snapshot-only restart, but the new regression is not authored in the repository's required result-test format and the exact default/override split-size values are not directly asserted.
  2. Scope: The production changes are small and focused on the job/provider lifecycle and the two MySQL/PostgreSQL reader construction paths.
  3. Concurrency: Job state is serialized by the existing job write lock, while provider split/progress state remains protected by splitsLock. The new transient persistence timestamp is only used inside the job-locked commit path. Async split enumeration keeps its existing publish recheck, no provider RPC was moved under splitsLock, no new nested lock order was introduced, and I found no deadlock or heavyweight-operation regression.
  4. Lifecycle: Normal execution, recovery scheduling, EditLog replay, checkpoint-image restore, transaction callbacks, and finished-job cleanup were traced. A crash before the terminal job log resumes from conservative durable offset state and rediscovers the end; a crash after it replays status and finish time. Callback registration/removal remains balanced. There are no C++ static-initialization concerns in this FE/Java-only change.
  5. Configuration: The persistence interval is mutable and read on every eligible snapshot commit, so changes take effect without restart; a new master's transient zero timestamp forces an immediate first journal. Both reader paths apply 40960 only when snapshot_split_size is absent, while validated explicit values still win.
  6. Compatibility: No RPC, function-symbol, or storage-format contract changes were introduced. The timer is transient, old serialized jobs safely use its zero default, and existing provider JSON fields remain readable across replay/rolling upgrade.
  7. Parallel paths: Base and TVF providers, normal and recovery completion, cloud and non-cloud scheduling, MySQL and PostgreSQL readers, and all startup modes were checked. The corresponding paths are updated consistently.
  8. Conditional checks: The non-empty binlog-offset invariant is enforced before snapshot-state cleanup in both providers and is the direct precondition for entering durable binlog state. Broader source-specific validation of malformed non-empty ALTER offsets is pre-existing behavior, not a regression introduced by this cleanup.
  9. Test coverage: The new tests exercise the core timing/replay/cleanup cases and a snapshot-only restart. The regression must be revised to hardcode its ordinary single-table name and expose deterministic status/count checks through named qt/order_qt output rather than raw Groovy assertions. A separate narrow gap is the lack of a direct reader assertion for absent property -> 40960 and explicit property -> exact configured value (including both source types); I treated that as nonblocking because the implementation is a direct validated getOrDefault path. Per the review environment's constraints, I performed static review only and did not run builds or tests.
  10. Test results: The deterministic expectations are consistent with the traced behavior, but the suite uses def tableName and raw asserts for stable query results, contrary to the mandatory regression-test rules, and the matching runner-generated .out file is absent. That result artifact must be generated through the prescribed regression runner, not handwritten.
  11. Observability: Existing job/offset logs retain identifiers and error context sufficient for recovery diagnosis. The new rate-limited path removes repeated persistence work and does not need a per-commit INFO log or a new critical-path metric.
  12. Transactions, persistence, and data writes: Offset update/replay ordering, EditLog coverage, image restore, and transaction callback behavior remain consistent. A failover during a suppressed snapshot interval can move the durable resume point backward and reprocess work, but cannot advance it past unjournaled data or skip a source range; that is consistent with the existing at-least-once/unique-key target path.
  13. FE-BE variables: No new variable or protocol field needs propagation between FE and BE.
  14. Performance and other issues: The change removes redundant snapshot EditLog writes and uses a larger default split size to reduce split overhead. Remaining per-commit serialization is pre-existing; I found no new CPU, memory, allocation, or algorithmic issue and no other blocking defect.

There was no additional review-focus guidance, so this conclusion covers the whole PR. Please fix the regression-test authoring issue before merge.

SELECT COUNT(*), COUNT(DISTINCT id)
FROM ${currentDb}.${tableName}
"""
assert rowsBeforeRestart.size() == 1

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Please bring this new regression into compliance with the repository's mandatory test rules. For a single ordinary table, hardcode snapshot_finished_restart_fe instead of def tableName; and record these deterministic pre/post-restart count/status checks through named qt/order_qt queries with the runner-generated .out, rather than raw Groovy asserts. The Awaitility readiness and dynamic job-ID checks can remain operational assertions.

@hello-stephen

Copy link
Copy Markdown
Contributor

FE Regression Coverage Report

Increment line coverage 4.88% (2/41) 🎉
Increment coverage report
Complete coverage report

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants